Guía visual de módulos y funciones

Documentación operativa de gss_bi_udfs

Esta versión está pensada para navegar rápido por módulos, pero sin perder la referencia detallada. Cada función tiene un bloque principal de lectura breve y un enlace desplegable para ver la firma, parámetros, retorno y ejemplo de uso al estilo de la documentación original.

Mapa del paquete

  • utils: Spark, ambiente, catálogos y metadata.
  • io: lectura de archivos, vistas temporales y escritura Delta.
  • transforms: hash técnico y registro default.
  • merges: histórico SCD2 sobre Delta.

Navegación sugerida

  • Usar el menú lateral para ubicar el módulo.
  • Leer el bloque “Para qué sirve”.
  • Abrir “Ver referencia completa” solo en las funciones que interesen.
Vista general

Resumen ejecutivo de módulos

La idea es identificar rápido qué problema resuelve cada módulo antes de entrar a la referencia fina de funciones.

utils

Centraliza el contexto técnico del paquete: SparkSession, ambiente, nombres de catálogo y metadata de tablas.

io

Es la capa de entrada y salida. Resuelve cargas raw, parquets, vistas temporales y persistencia Delta.

transforms

Aporta helpers de enriquecimiento de DataFrames, especialmente para claves técnicas y filas default.

merges

Especializado en versionado histórico SCD2, con cierre de vigencias e inserción de nuevas versiones.

Configuración y contexto

Módulo utils

Este módulo resuelve dependencias transversales del paquete: acceso a Spark, ambiente, catálogos y ubicación de tablas.

Para qué sirve

  • Evita repetir la lógica de inicialización de Spark.
  • Normaliza catálogos y rutas según ENV.
  • Expone metadata clave para lectura y escritura segura.
Spark Configuración

set_spark

set_spark(spark)

Registra una sesión Spark global para reutilizarla desde el resto del paquete.

Ver referencia completa

Configura una SparkSession global para uso interno del módulo.

ParámetroTipoDescripción
sparkSparkSessionSesión Spark a registrar globalmente.
Retorna: None.
from pyspark.sql import SparkSession
from gss_bi_udfs import utils

spark = SparkSession.builder.getOrCreate()
utils.set_spark(spark)
Spark Fallbacks

get_spark

get_spark(spark=None)

Devuelve la sesión Spark disponible usando una cadena ordenada de fallbacks.

Ver referencia completa

Obtiene Spark en este orden: parámetro explícito, sesión configurada por set_spark, builtins.spark o builder.getOrCreate().

ParámetroTipoDescripción
sparkSparkSession | NoneSesión opcional explícita.
Retorna: SparkSession activa.
spark = utils.get_spark()
Ambiente Helper

get_env

get_env(default="dev")

Lee el entorno de ejecución desde la variable ENV.

Ver referencia completa

Obtiene el entorno de ejecución y usa un default cuando la variable no está definida.

ParámetroTipoDescripción
defaultstrValor por defecto si ENV no existe.
Retorna: str con el entorno, por ejemplo dev, qa o pro.
env = utils.get_env(default="dev")
Catálogos Ambiente

get_env_catalog

get_env_catalog(catalog)

Ajusta el catálogo según el ambiente actual.

Ver referencia completa

Si ENV=pro devuelve el catálogo base. En otros casos agrega el sufijo del ambiente.

ParámetroTipoDescripción
catalogstrNombre base del catálogo.
Retorna: str con catálogo ajustado por ambiente.
catalogo = utils.get_env_catalog("fi_comunes")
Tablas Naming

get_env_table_path

get_env_table_path(catalog, table_path)

Compone el identificador final de tabla usando catálogo ajustado y ruta lógica.

Ver referencia completa

Concatena el catálogo modificado por ambiente con la ruta schema.tabla.

ParámetroTipoDescripción
catalogstrCatálogo base.
table_pathstrRuta esquema.tabla.
Retorna: str con la tabla completa.
path = utils.get_env_table_path("fi_comunes", "silver.dim_cliente")
Metadata Infraestructura

get_schema_root_location

get_schema_root_location(spark=None, catalog=None, schema=None)

Consulta la ubicación física real del esquema dentro del catálogo correspondiente.

Ver referencia completa

Ejecuta DESCRIBE SCHEMA EXTENDED y extrae el valor RootLocation.

ParámetroTipoDescripción
sparkSparkSession | NoneSesión opcional.
catalogstrCatálogo base.
schemastrEsquema dentro del catálogo.
Retorna: str con el path físico del esquema.
root = utils.get_schema_root_location(
    spark=spark,
    catalog="fi_comunes",
    schema="silver"
)
Metadata Escritura segura

get_table_info

get_table_info(spark=None, full_table_name=None, catalog=None, schema=None, table=None)

Resuelve nombre final, path, existencia y metadatos de una tabla.

Ver referencia completa

Devuelve información integral de tabla: nombre completo, path, existencia, provider y tipo.

ParámetroTipoDescripción
sparkSparkSession | NoneSesión opcional.
full_table_namestr | NoneFormato catalog.schema.table.
catalog / schema / tablestr | NoneAlternativa al parámetro anterior.
Retorna: dict con catalog, schema, table, full_table_name, path, exists, provider y table_type.
info = utils.get_table_info(
    spark=spark,
    full_table_name="fi_comunes.silver.dim_cliente"
)
print(info["full_table_name"], info["path"], info["exists"])
Defaults Tipos PySpark

get_default_value_by_type

get_default_value_by_type(dtype)

Mapea un tipo PySpark a un valor default técnico consistente.

Ver referencia completa

Devuelve una Column con valores como -999, N/A, 1900-01-01 o False según el tipo de dato.

ParámetroTipoDescripción
dtypepyspark.sql.types.DataTypeTipo de dato del campo.
Retorna: Column con valor por defecto.
default_col = utils.get_default_value_by_type(df.schema["fecha_alta"].dataType)
Lectura y persistencia

Módulo io

Este módulo concentra el flujo operativo de entrada y salida: archivos raw, parquets, vistas temporales y tablas Delta.

Bloques funcionales

  • Lectura del último archivo disponible por carpeta.
  • Carga de parquets versionados desde bronze.
  • Registro masivo de vistas temporales.
  • Persistencia estándar y por snapshot en Delta.
Raw files CSV/TXT

load_latest_delimited_file

load_latest_delimited_file(spark=None, source_file=None, env=None, delimiter=",", header=True, infer_schema=True, ...)

Carga el último archivo delimitado compatible dentro de la carpeta indicada.

Ver referencia completa

Está orientada a .csv, .txt y variantes similares leídas de forma nativa con Spark.

ParámetroTipoDescripción
sparkSparkSession | NoneSesión opcional.
source_filestrRuta relativa dentro de /Volumes/bronze/files/{env}/.
envstr | NoneAmbiente de lectura.
delimiterstrSeparador del archivo.
headerboolIndica si el archivo tiene encabezado.
infer_schemaboolSi Spark debe inferir tipos.
encodingstrEncoding del archivo.
quotestrCarácter quote.
escapestrCarácter escape.
multilineboolSoporte para registros multilinea.
extensionsiterable | NoneExtensiones admitidas.
Retorna: DataFrame de Spark o None si no encuentra archivos.
from gss_bi_udfs import io

df_clientes = io.load_latest_delimited_file(
    spark=spark,
    source_file="comercial/clientes_vigentes/periodo=2026-03",
    delimiter=";",
    header=True,
    encoding="latin1"
)
Excel Raw files

load_latest_excel_file

load_latest_excel_file(spark=None, source_file=None, env=None, header=True, sheet_name=0)

Carga el último archivo Excel y lo convierte a DataFrame Spark usando pandas.

Ver referencia completa

Busca el último .xls o .xlsx dentro de la carpeta solicitada.

ParámetroTipoDescripción
sparkSparkSession | NoneSesión opcional.
source_filestrRuta relativa del dataset.
envstr | NoneAmbiente de lectura.
headerboolSi la primera fila contiene encabezado.
sheet_namestr | intHoja a leer por nombre o índice.
Retorna: DataFrame de Spark o None si no hay archivos.
df_excel = io.load_latest_excel_file(
    spark=spark,
    source_file="comercial/clientes_vigentes/periodo=2026-03",
    header=True,
    sheet_name="Hoja1"
)
Dispatcher Formato único

load_latest_file

load_latest_file(spark=None, source_file=None, env=None, file_format=None, ...)

Elige internamente la estrategia de lectura adecuada según el formato declarado.

Ver referencia completa

Soporta excel, xls, xlsx, csv, txt, tsv y delimited.

ParámetroTipoDescripción
sparkSparkSession | NoneSesión opcional.
source_filestrRuta relativa del dataset versionado.
envstr | NoneAmbiente de lectura.
file_formatstrFormato del archivo a leer.
headerboolSi el archivo contiene encabezado.
delimiterstrSeparador para formatos delimitados.
infer_schemaboolInferencia de tipos.
encodingstrEncoding del origen.
quotestrCarácter quote.
escapestrCarácter escape.
multilineboolSoporte para multilinea.
sheet_namestr | intHoja en caso de Excel.
Retorna: DataFrame de Spark.
df = io.load_latest_file(
    spark=spark,
    source_file="comercial/clientes_vigentes/periodo=2026-03",
    file_format="delimited",
    delimiter="|",
    header=False,
    infer_schema=False
)
Bronze Parquet

load_latest_parquet

load_latest_parquet(spark=None, data_base=None, schema=None, table=None, env=None)

Recupera el parquet más reciente de una tabla en bronze y lo devuelve como DataFrame.

Ver referencia completa

Carga el parquet más reciente desde /Volumes/bronze/{db}_{schema}/{env}/{table}/.

ParámetroTipoDescripción
sparkSparkSession | NoneSesión opcional.
data_basestrNombre de base lógica.
schemastrNombre de esquema.
tablestrNombre de tabla.
envstr | NoneAmbiente de lectura.
Retorna: DataFrame de Spark o None.
df_clientes = io.load_latest_parquet(
    spark=spark,
    data_base="timepro",
    schema="insudb",
    table="clientes",
    env="dev"
)
Temp views Retorna diccionario

return_parquets_and_register_temp_views

return_parquets_and_register_temp_views(spark=None, tables_load=None, verbose=False, env=None)

Carga varios parquets, materializa vistas temporales y devuelve los DataFrames en un diccionario.

Ver referencia completa

Permite definir una configuración de tablas y vistas para levantar varias fuentes de una sola vez.

ParámetroTipoDescripción
sparkSparkSession | NoneSesión opcional.
tables_loaddictMapa de carga por base, esquema, tabla y vista.
verboseboolImprime mensajes de estado.
envstr | NoneAmbiente de lectura.
Retorna: dict con claves db.schema.table y valores DataFrame.
tables_load = {
    "timepro": {
        "insudb": [
            {"table": "clientes", "view": "vw_clientes"}
        ]
    }
}

dfs = io.return_parquets_and_register_temp_views(
    spark=spark,
    tables_load=tables_load,
    verbose=True
)
Temp views Sin retorno

parquets_register_temp_views

parquets_register_temp_views(spark=None, tables_load=None, verbose=False, env=None)

Versión orientada solo a materializar vistas temporales desde parquets.

Ver referencia completa

Misma lógica general de carga que la función anterior, pero sin retornar el diccionario.

ParámetroTipoDescripción
sparkSparkSession | NoneSesión opcional.
tables_loaddictConfiguración de tablas y vistas.
verboseboolMensajes de materialización.
envstr | NoneAmbiente de lectura.
Retorna: None.
io.parquets_register_temp_views(
    spark=spark,
    tables_load=tables_load,
    verbose=True
)
Excel Compatibilidad

load_latest_excel

load_latest_excel(spark=None, source_file=None, env=None)

Función simple de compatibilidad para cargas Excel con encabezado activo.

Ver referencia completa

Internamente delega a load_latest_excel_file y retorna None si falla.

ParámetroTipoDescripción
sparkSparkSession | NoneSesión opcional.
source_filestrRuta relativa del archivo dentro de bronze/files.
envstr | NoneAmbiente de lectura.
Retorna: DataFrame de Spark o None.
df_legacy = io.load_latest_excel(
    spark=spark,
    source_file="comercial/clientes_vigentes/periodo=2026-03",
    env="dev"
)
Excel batch Retorna diccionario

return_excels_and_register_temp_views

return_excels_and_register_temp_views(spark=None, files_load=None, verbose=False, env=None)

Carga un conjunto de Excels y registra vistas temporales, devolviendo además los DataFrames cargados.

Ver referencia completa

Es útil cuando un proceso trabaja con varios archivos por dominio y subdominio.

ParámetroTipoDescripción
sparkSparkSession | NoneSesión opcional.
files_loaddictConfiguración dominio/subdominio/lista de archivos y vistas.
verboseboolMensajes de materialización.
envstr | NoneAmbiente de lectura.
Retorna: dict con claves dominio.subdominio.archivo y valores DataFrame.
files_load = {
    "comercial": {
        "clientes": [
            {"file": "padron", "view": "vw_padron"}
        ]
    }
}

dfs = io.return_excels_and_register_temp_views(
    spark=spark,
    files_load=files_load,
    verbose=True
)
Excel batch Solo vistas

excels_register_temp_views

excels_register_temp_views(spark=None, files_load=None, verbose=False, env=None)

Materializa vistas temporales desde múltiples archivos Excel sin devolver estructura adicional.

Ver referencia completa

Es la variante de ejecución más directa cuando solo importa dejar vistas registradas en Spark.

ParámetroTipoDescripción
sparkSparkSession | NoneSesión opcional.
files_loaddictConfiguración de archivos y vistas.
verboseboolMensajes de materialización.
envstr | NoneAmbiente de lectura.
Retorna: None.
io.excels_register_temp_views(
    spark=spark,
    files_load=files_load,
    verbose=True
)
Dispatcher Acciones

load_and_materialize_views

load_and_materialize_views(action, **kwargs)

Ejecuta una estrategia de carga a partir del nombre de una acción predefinida.

Ver referencia completa

Permite delegar la ejecución a funciones como return_parquets_and_register_temp_views o excels_register_temp_views.

ParámetroTipoDescripción
actionstrNombre de la acción a ejecutar.
kwargsdictArgumentos a pasar a la función seleccionada.
Retorna: dict o resultado de la acción ejecutada; si la acción no existe, retorna un dict vacío.
result = io.load_and_materialize_views(
    "return_parquets_and_register_temp_views",
    spark=spark,
    tables_load=tables_load,
    verbose=True
)
Delta Persistencia base

save_table_to_delta

save_table_to_delta(spark=None, df=None, catalog=None, schema=None, table_name=None, mode="overwrite")

Guarda un DataFrame en Delta usando nombre lógico y path físico resueltos automáticamente.

Ver referencia completa

Usa get_table_info para determinar la ubicación y el nombre completo de la tabla destino.

ParámetroTipoDescripción
sparkSparkSession | NoneSesión activa de Spark.
dfDataFrameDataFrame a persistir.
catalogstrCatálogo destino.
schemastrEsquema destino.
table_namestrNombre de la tabla destino.
modestrModo de escritura Spark, por ejemplo overwrite o append.
Retorna: None.
io.save_table_to_delta(
    spark=spark,
    df=df_resultado,
    catalog="fi_comunes",
    schema="silver",
    table_name="dim_cliente",
    mode="overwrite"
)
Delta Reemplazo por rango

save_snapshot_to_delta

save_snapshot_to_delta(spark=None, df=None, catalog=None, schema=None, table_name=None, snapshot_column=None)

Reemplaza el rango existente de un snapshot y luego escribe los nuevos registros en Delta.

Ver referencia completa

Está pensada para tablas fact cargadas por período o time_key, validando el rango presente en el DataFrame origen.

ParámetroTipoDescripción
sparkSparkSession | NoneSesión activa de Spark.
dfDataFrameDataFrame a guardar.
catalogstrCatálogo destino.
schemastrEsquema destino.
table_namestrTabla destino.
snapshot_columnstrColumna que define el período del snapshot.
Retorna: None.
io.save_snapshot_to_delta(
    spark=spark,
    df=df_fact,
    catalog="fi_comunes",
    schema="gold",
    table_name="fact_ventas",
    snapshot_column="periodo"
)
Lectura rápida del módulo: si la pregunta es “cómo leo esto”, “cómo dejo vistas temporales” o “cómo persisto en Delta”, la función seguramente está en io.
Enriquecimiento de DataFrames

Módulo transforms

Módulo corto y puntual. Su foco está en agregar claves técnicas y construir registros default a partir del esquema.

Para qué sirve

  • Generar identificadores hash repetibles.
  • Construir filas “unknown” o “N/A” manteniendo tipos correctos.
Hash PK técnica

add_hashid

add_hashid(df, columns, new_col_name="hashid")

Concatena columnas como string, calcula xxhash64 y deja la nueva columna al inicio del DataFrame.

Ver referencia completa

Sirve para crear una PK técnica rápida y determinística a partir de columnas de negocio.

ParámetroTipoDescripción
dfDataFrameDataFrame de Spark de entrada.
columnslist[str]Lista de columnas a concatenar para el hash.
new_col_namestrNombre de la nueva columna hash.
Retorna: DataFrame con la columna hash agregada.
from gss_bi_udfs.transforms import add_hashid

df_hash = add_hashid(
    df=df_clientes,
    columns=["documento", "sucursal"],
    new_col_name="cliente_hash"
)
Defaults Dimensiones

get_default_record

get_default_record(spark, df)

Construye un DataFrame con una sola fila de valores default respetando el esquema del DataFrame origen.

Ver referencia completa

Usa get_default_value_by_type para asignar un valor técnico por cada tipo del esquema.

ParámetroTipoDescripción
sparkSparkSessionSesión activa de Spark.
dfDataFrameDataFrame del cual se toma el esquema.
Retorna: DataFrame con un único registro default.
from gss_bi_udfs.transforms import get_default_record

df_default = get_default_record(spark, df_dim_cliente)
Versionado histórico

Módulo merges

Este módulo está orientado a mantener dimensiones históricas con SCD Tipo 2 sobre Delta Lake.

Qué resuelve

  • Carga inicial si la tabla destino no existe.
  • Cierre de la versión activa cuando cambian atributos.
  • Inserción de nuevas versiones con nuevas vigencias.
  • Baja lógica de registros que ya no están en el origen activo.
SCD2 Histórico Modelo consistente

merge_scd2

merge_scd2(spark=None, df_dim_src=None, table_name=None, business_keys=None, surrogate_key=None, eow_date="9999-12-31")

Implementa Slowly Changing Dimension Tipo 2 usando business keys, detección de cambios y control de vigencias con valid_from y valid_to.

Ver referencia completa

La surrogate key se genera internamente como un hash de business_keys + valid_from. La función compara todas las columnas del origen excepto business keys y columnas de vigencia.

ParámetroTipoDescripción
sparkSparkSession | NoneSesión activa de Spark.
df_dim_srcDataFrameDataFrame con la dimensión actual a nivel de negocio.
table_namestrNombre completo de la tabla Delta destino.
business_keysstr | list[str]Clave simple o compuesta que identifica la entidad de negocio.
surrogate_keystrNombre de la clave primaria física.
eow_datestrFecha de fin de vigencia para registros activos.
Retorna: None. Ejecuta el merge directamente sobre la tabla Delta.
from gss_bi_udfs.merges import merge_scd2

merge_scd2(
    spark=spark,
    df_dim_src=df_clientes,
    table_name="fi_comunes.silver.dim_cliente",
    business_keys=["nro_documento"],
    surrogate_key="cliente_sk"
)
Importante: el DataFrame origen no debería incluir columnas volátiles como timestamps de carga o identificadores técnicos transitorios, porque serían interpretadas como cambios reales y generarían nuevas versiones.
Cómo combinar módulos

Flujo recomendado de uso

Este ejemplo resume cómo se conectan los módulos en un caso BI típico: leer, transformar y persistir.

from pyspark.sql import SparkSession
from gss_bi_udfs import io, transforms, utils
from gss_bi_udfs.merges import merge_scd2

spark = SparkSession.builder.getOrCreate()
utils.set_spark(spark)

df_raw = io.load_latest_file(
    source_file="ventas/facturacion/clientes",
    file_format="csv",
    delimiter=";",
    header=True
)

df_dim = transforms.add_hashid(df_raw, ["documento"], "cliente_hash")

merge_scd2(
    df_dim_src=df_dim,
    table_name="fi_comunes.silver.dim_cliente",
    business_keys=["documento"],
    surrogate_key="cliente_sk"
)